GuestOps/src/GuestOps.Api/GoogleMailbox.cs

172 lines
15 KiB
C#

using System.Text;
using System.Text.Json;
using System.Net;
using System.Security.Cryptography;
using Microsoft.AspNetCore.DataProtection;
namespace GuestOps.Web;
public sealed class GoogleFailure(string kind,HttpStatusCode status,TimeSpan? retryAfter=null):Exception("Google request failed: "+kind)
{
public string Kind {get;}=kind;
public HttpStatusCode Status {get;}=status;
public TimeSpan? RetryAfter {get;}=retryAfter;
}
public sealed class GoogleMailbox(HttpClient http, IConfiguration config, IStore store, IDataProtectionProvider protection)
{
private readonly IDataProtector protector = protection.CreateProtector("GoogleMailbox.refresh.v1");
private string ClientId => config["Google:ClientId"] ?? "";
private string ClientSecret => config["Google:ClientSecret"] ?? "";
private string Callback => (config["PublicUrl"] ?? "https://sandbox-guestops.futuresens.co.uk").TrimEnd('/') + "/api/integrations/google/callback";
public bool Configured => ClientId.Length > 0 && ClientSecret.Length > 0;
public bool SendingConfigured => Configured && config.GetValue<bool>("Google:EnableSending");
public string AuthorizationUrl(string state) => "https://accounts.google.com/o/oauth2/v2/auth?" + string.Join("&", new Dictionary<string,string> {
["client_id"] = ClientId, ["redirect_uri"] = Callback, ["response_type"] = "code", ["scope"] = "https://www.googleapis.com/auth/gmail.readonly" + (SendingConfigured ? " https://www.googleapis.com/auth/gmail.send" : ""), ["access_type"] = "offline", ["prompt"] = "consent", ["state"] = state
}.Select(x => Uri.EscapeDataString(x.Key) + "=" + Uri.EscapeDataString(x.Value)));
async Task<JsonElement> Token(Dictionary<string, string> data, CancellationToken ct = default)
{
data["client_id"] = ClientId; data["client_secret"] = ClientSecret;
using var response = await http.PostAsync("https://oauth2.googleapis.com/token", new FormUrlEncodedContent(data), ct);
await Check(response,true,ct);
return JsonDocument.Parse(await response.Content.ReadAsStringAsync(ct)).RootElement.Clone();
}
async Task<JsonElement> Read(string path, string token, CancellationToken ct = default)
{
using var request = new HttpRequestMessage(HttpMethod.Get, "https://gmail.googleapis.com/gmail/v1/users/me/" + path);
request.Headers.Authorization = new("Bearer", token);
using var response = await http.SendAsync(request, ct); await Check(response,false,ct);
return JsonDocument.Parse(await response.Content.ReadAsStringAsync(ct)).RootElement.Clone();
}
static async Task Check(HttpResponseMessage response,bool token,CancellationToken ct)
{
if(response.IsSuccessStatusCode)return;
var reason="";try{using var json=JsonDocument.Parse(await response.Content.ReadAsStringAsync(ct));if(json.RootElement.TryGetProperty("error",out var error)){if(error.ValueKind==JsonValueKind.String)reason=error.GetString()??"";else if(error.TryGetProperty("errors",out var errors)&&errors.GetArrayLength()>0&&errors[0].TryGetProperty("reason",out var value))reason=value.GetString()??"";}}catch(JsonException){}
var kind=token?(reason=="invalid_grant"?"ReconnectRequired":reason=="invalid_client"||response.StatusCode==HttpStatusCode.Unauthorized?"Configuration":"Temporary"):
response.StatusCode==HttpStatusCode.Unauthorized?"ReconnectRequired":response.StatusCode==HttpStatusCode.NotFound?"NotFound":response.StatusCode==HttpStatusCode.BadRequest?"BadRequest":response.StatusCode==HttpStatusCode.Forbidden&&reason is not ("rateLimitExceeded" or "userRateLimitExceeded")?"AccessDenied":"Temporary";
var delay=response.Headers.RetryAfter?.Delta??(response.Headers.RetryAfter?.Date-DateTimeOffset.UtcNow);
throw new GoogleFailure(kind,response.StatusCode,delay);
}
public async Task Connect(string hotel, string code,DateTime? startedAt=null,string expectedEmail="")
{
if (string.IsNullOrWhiteSpace(code)) throw new InvalidOperationException("No authorization code.");
var tokens = await Token(new() { ["code"] = code, ["redirect_uri"] = Callback, ["grant_type"] = "authorization_code" });
var profile = await Read("profile", tokens.GetProperty("access_token").GetString()!);
var email = profile.GetProperty("emailAddress").GetString()!.ToLowerInvariant();
if(!Input.Email(email)||expectedEmail.Length>0&&email!=expectedEmail)throw new InvalidOperationException("Choose the expected Google mailbox.");
var prior = (await store.List<Mailbox>(hotel)).SingleOrDefault(x => x.Email == email);
if(prior?.ConnectionChangedAt>startedAt.GetValueOrDefault(DateTime.MinValue))throw new MailboxConflict();
var mailbox = prior ?? new Mailbox { HotelId = hotel, Email = email };
var version=mailbox.Version;
// No token reuse across hotels. Mongo's unique mailbox-email index prevents
// accidental connection of one shared mailbox to two hotel workspaces.
mailbox.ProtectedRefreshToken = protector.Protect(tokens.GetProperty("refresh_token").GetString()!);
mailbox.CanSend = tokens.TryGetProperty("scope", out var scopes) && scopes.GetString()!.Split(' ').Contains("https://www.googleapis.com/auth/gmail.send");
if(string.IsNullOrWhiteSpace(tokens.GetProperty("refresh_token").GetString()))throw new InvalidOperationException("Google did not return offline access.");
if(!tokens.TryGetProperty("scope",out var granted)||!granted.GetString()!.Split(' ').Contains("https://www.googleapis.com/auth/gmail.readonly"))throw new InvalidOperationException("Read permission is required.");
mailbox.Status = "Connected"; mailbox.SyncError = "";mailbox.SyncErrorCode="";mailbox.NextAttemptAt=null;mailbox.FailureCount=0;mailbox.PageToken="";
mailbox.ConnectionEpoch=Guid.NewGuid().ToString("N");mailbox.ConnectionChangedAt=DateTime.UtcNow;mailbox.Version++;
if(!(prior==null?await store.TryInsertMailbox(mailbox):await store.Replace(hotel,mailbox.Id,version,mailbox)))throw new MailboxConflict();
}
public async Task Sync(Mailbox mailbox, CancellationToken ct)
{
if(mailbox.NextAttemptAt>DateTime.UtcNow||!await MailboxManagement.Current(store,mailbox))return;
mailbox.LastAttemptAt=DateTime.UtcNow;await store.SaveSync(mailbox);
var token = await AccessToken(mailbox,ct);
var after = new DateTimeOffset(mailbox.WindowStart).ToUnixTimeSeconds();
var before = new DateTimeOffset(mailbox.WindowEnd).ToUnixTimeSeconds();
var query = Uri.EscapeDataString($"in:inbox after:{after} before:{before}");
var path = "messages?maxResults=25&q=" + query + (mailbox.PageToken.Length > 0 ? "&pageToken=" + Uri.EscapeDataString(mailbox.PageToken) : "");
JsonElement page;
try{page=await Read(path, token, ct);}
catch(GoogleFailure ex) when(ex.Status==HttpStatusCode.BadRequest&&mailbox.PageToken.Length>0)
{
mailbox.PageToken="";mailbox.SyncErrorCode="CheckpointRestart";mailbox.SyncError="Google rejected the saved page. The worker will restart this import window without duplicating messages.";mailbox.NextAttemptAt=DateTime.UtcNow.AddMinutes(1);await store.SaveSync(mailbox);return;
}
if (page.TryGetProperty("messages", out var items)) foreach (var item in items.EnumerateArray())
{
if(!await MailboxManagement.Current(store,mailbox))return;
var id = item.GetProperty("id").GetString()!;
JsonElement message;try{message=await Read("messages/" + Uri.EscapeDataString(id) + "?format=full", token, ct);}catch(GoogleFailure ex) when(ex.Status==HttpStatusCode.NotFound){continue;}
var payload = message.GetProperty("payload");
string Header(string name) => payload.GetProperty("headers").EnumerateArray().Where(x => string.Equals(x.GetProperty("name").GetString(), name, StringComparison.OrdinalIgnoreCase)).Select(x => x.GetProperty("value").GetString()).FirstOrDefault() ?? "";
var auto = Header("Auto-Submitted");
if ((auto.Length > 0 && auto != "no") || Header("List-Id").Length > 0 || Header("Return-Path").Trim() == "<>") continue;
var body = PlainText(payload);
var row = new Conversation { HotelId = mailbox.HotelId, MailboxId = mailbox.Id, ProviderMessageId = id, ProviderThreadId = message.GetProperty("threadId").GetString()!, From = Header("From"), Subject = Header("Subject"), Body = body.Length > 0 ? body[..Math.Min(body.Length, 30000)] : "This message has no plain-text body. Open it in Gmail to read it.", ReceivedAt = DateTimeOffset.FromUnixTimeMilliseconds(long.Parse(message.GetProperty("internalDate").GetString()!)).UtcDateTime };
row.ReplyAddress = ReplyMime.Address(Header("Reply-To").Length > 0 ? Header("Reply-To") : Header("From"));
row.RfcMessageId = Header("Message-ID");
row.AutoReplyHeadersEligible = FaqMatcher.HeadersEligible(payload,mailbox.Email);
if(!await MailboxManagement.Current(store,mailbox))return;
await store.Import(row); // deduplicated before advancing the page checkpoint
}
mailbox.PageToken = page.TryGetProperty("nextPageToken", out var next) ? next.GetString()! : "";
if (mailbox.PageToken.Length == 0) { mailbox.WindowStart = mailbox.WindowEnd.AddMinutes(-5); mailbox.WindowEnd = DateTime.UtcNow; }
mailbox.LastSyncAt = DateTime.UtcNow; mailbox.SyncError = "";mailbox.SyncErrorCode="";mailbox.FailureCount=0;mailbox.NextAttemptAt=null;
await store.SaveSync(mailbox);
}
static string PlainText(JsonElement payload)
{
if (payload.TryGetProperty("mimeType", out var mime) && mime.GetString() == "text/plain" && payload.TryGetProperty("body", out var body) && body.TryGetProperty("data", out var data))
{
var raw = data.GetString()!.Replace('-', '+').Replace('_', '/'); raw = raw.PadRight((raw.Length + 3) / 4 * 4, '=');
return Encoding.UTF8.GetString(Convert.FromBase64String(raw));
}
return payload.TryGetProperty("parts", out var parts) ? string.Join("\n", parts.EnumerateArray().Select(PlainText).Where(x => x.Length > 0)) : "";
}
public async Task<string> AccessToken(Mailbox mailbox, CancellationToken ct)
{
if(!await MailboxManagement.Current(store,mailbox))throw new MailboxConflict();
try
{
var token = await Token(new() { ["refresh_token"] = protector.Unprotect(mailbox.ProtectedRefreshToken), ["grant_type"] = "refresh_token" }, ct);
return token.GetProperty("access_token").GetString()!;
}
catch(Exception ex) when(ex is CryptographicException || ex is GoogleFailure {Kind:"ReconnectRequired"}) {await RecordFailure(mailbox,ex);throw;}
}
public async Task RecordFailure(Mailbox mailbox,Exception ex)
{
if(ex is MailboxConflict)return;
mailbox.LastAttemptAt=DateTime.UtcNow;mailbox.FailureCount=Math.Min(mailbox.FailureCount+1,20);
var kind=ex is GoogleFailure failure?failure.Kind:ex is CryptographicException?"ReconnectRequired":"Temporary";
mailbox.SyncErrorCode=kind;
if(kind=="ReconnectRequired") {mailbox.Status="NeedsReconnect";mailbox.NextAttemptAt=null;mailbox.SyncError="Google access is unavailable. Ask the hotel owner to reconnect this mailbox.";}
else
{
var seconds=kind is "Configuration" or "AccessDenied"?3600:Math.Min(3600,60*Math.Pow(2,mailbox.FailureCount-1));
if(ex is GoogleFailure g&&g.RetryAfter is {} delay)seconds=Math.Max(seconds,Math.Min(86400,delay.TotalSeconds));
mailbox.NextAttemptAt=DateTime.UtcNow.AddSeconds(seconds+Random.Shared.Next(0,31));
mailbox.SyncError=kind is "Configuration" or "AccessDenied"?"Google rejected the app configuration or permissions. Ask the administrator to review access. A later check is scheduled.":"Google synchronization is temporarily unavailable. The worker will retry automatically.";
}
await store.SaveSync(mailbox);
}
public async Task<bool> AutoReplyThreadUnchanged(Conversation message,string token,CancellationToken ct)
{
var thread=await Read("threads/"+Uri.EscapeDataString(message.ProviderThreadId)+"?format=metadata",token,ct);
if(!thread.TryGetProperty("messages",out var messages)||messages.GetArrayLength()!=1)return false;
var only=messages[0];return only.GetProperty("id").GetString()==message.ProviderMessageId&&only.TryGetProperty("labelIds",out var labels)&&labels.EnumerateArray().Any(x=>x.GetString()=="INBOX")&&!labels.EnumerateArray().Any(x=>x.GetString()=="SENT");
}
public async Task<string> Send(Conversation message, string token, string raw, CancellationToken ct)
{
using var request = new HttpRequestMessage(HttpMethod.Post, "https://gmail.googleapis.com/gmail/v1/users/me/messages/send") { Content = JsonContent.Create(new { raw, threadId = message.ProviderThreadId }) };
request.Headers.Authorization = new("Bearer", token);
using var response = await http.SendAsync(request, ct);
// Once SendAsync is entered, any exception or unexpected response is ambiguous.
// This adapter never automatically retries a provider send.
response.EnsureSuccessStatusCode();
using var json = JsonDocument.Parse(await response.Content.ReadAsStringAsync(ct));
var id = json.RootElement.GetProperty("id").GetString();
return !string.IsNullOrWhiteSpace(id) ? id : throw new InvalidOperationException("No delivery ID returned.");
}
public async Task<string?> FindSent(Conversation message, Mailbox mailbox, CancellationToken ct)
{
var token = await AccessToken(mailbox, ct);
var query = Uri.EscapeDataString("in:sent rfc822msgid:" + message.Delivery!.MessageId);
var page = await Read("messages?maxResults=2&q=" + query, token, ct);
if (!page.TryGetProperty("messages", out var items) || items.GetArrayLength() != 1) return null;
var id = items[0].GetProperty("id").GetString()!;
var found = await Read("messages/" + Uri.EscapeDataString(id) + "?format=metadata", token, ct);
var headers = found.GetProperty("payload").GetProperty("headers").EnumerateArray().ToArray();
string Header(string name) => headers.FirstOrDefault(h => h.GetProperty("name").GetString()!.Equals(name, StringComparison.OrdinalIgnoreCase)) is var h && h.ValueKind != JsonValueKind.Undefined ? h.GetProperty("value").GetString()! : "";
return found.GetProperty("labelIds").EnumerateArray().Any(x => x.GetString() == "SENT") && Header("Message-ID") == message.Delivery.MessageId && ReplyMime.Address(Header("To")) == message.Delivery.Recipient && ReplyMime.Address(Header("From")) == mailbox.Email ? id : null;
}
}