Skip to content
Open
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
10 changes: 8 additions & 2 deletions src/typeagent/emails/email_import.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
# Copyright (c) Microsoft Corporation.
# Licensed under the MIT License.

from datetime import datetime
from datetime import datetime, timezone
from email import message_from_string
from email.header import decode_header, Header, make_header
from email.message import Message
Expand All @@ -10,6 +10,7 @@
import re
from typing import Iterable, overload

from ..knowpro.universal_message import format_timestamp_utc
from .email_message import EmailMessage, EmailMessageMeta


Expand Down Expand Up @@ -100,7 +101,12 @@ def import_email_message(msg: Message, max_chunk_length: int) -> EmailMessage:
timestamp: str | None = None
timestamp_date = msg.get("Date", None)
if timestamp_date is not None:
timestamp = parsedate_to_datetime(timestamp_date).isoformat()
parsed_date = parsedate_to_datetime(timestamp_date)
if parsed_date.tzinfo is None:
# RFC 5322 "-0000": time is UTC but the origin zone is unknown.
parsed_date = parsed_date.replace(tzinfo=timezone.utc)
# Normalize to UTC so timestamps sort lexicographically across senders.
timestamp = format_timestamp_utc(parsed_date)
Comment thread
bmerkle marked this conversation as resolved.

# Get email body.
# If the email was a reply, then ensure we only pick up the latest response
Expand Down
25 changes: 18 additions & 7 deletions src/typeagent/storage/memory/timestampindex.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,9 @@
# Timestamp-to-text-range in-memory index (pre-SQLite prep).
#
# Contract (stable regardless of backing store):
# - add_timestamp(s) accepts ISO 8601 timestamps that are lexicographically sortable
# (Datetime.isoformat). Missing/None timestamps are ignored.
# - add_timestamp(s) accepts ISO 8601 timestamps with any UTC offset (naive
# timestamps are assumed to be UTC). They are normalized to UTC ("...Z") so that
# lexicographic order equals chronological order. Missing/None timestamps are ignored.
# - lookup_range(DateRange) returns items whose ISO timestamp t satisfies
# start <= t < end (end is exclusive). If end is None, treat as a point
# query with end = start + epsilon.
Expand All @@ -20,6 +21,7 @@

import bisect
from collections.abc import AsyncIterable, Callable
from datetime import timezone
from typing import Any

from ...knowpro.interfaces import (
Expand Down Expand Up @@ -56,8 +58,10 @@ async def lookup_range(self, date_range: DateRange) -> list[TimestampedTextRange
return self._lookup_range(date_range)

def _lookup_range(self, date_range: DateRange) -> list[TimestampedTextRange]:
start_at = date_range.start.isoformat()
stop_at = None if date_range.end is None else date_range.end.isoformat()
start_at = _normalize_datetime(date_range.start)
stop_at = (
None if date_range.end is None else _normalize_datetime(date_range.end)
)
return get_in_range(
self._ranges,
start_at,
Expand Down Expand Up @@ -101,11 +105,10 @@ def _insert_timestamp(
) -> bool:
if not timestamp:
return False
timestamp_datetime = Datetime.fromisoformat(timestamp)
entry: TimestampedTextRange = TimestampedTextRange(
range=text_range_from_message_chunk(message_ordinal),
# This string is formatted to be lexically sortable.
timestamp=timestamp_datetime.isoformat(),
# Normalized to UTC so that the string is lexically sortable.
timestamp=_normalize_datetime(Datetime.fromisoformat(timestamp)),
)
if in_order:
where = bisect.bisect_left(
Expand All @@ -117,6 +120,14 @@ def _insert_timestamp(
return True


def _normalize_datetime(dt: Datetime) -> str:
# Render as UTC with a fixed format ("Z" suffix, always microseconds) so that
# lexicographic order equals chronological order. Naive datetimes are UTC.
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc).isoformat(timespec="microseconds")[:-6] + "Z"


def get_in_range[T, S: Any](
values: list[T],
start_at: S,
Expand Down
10 changes: 10 additions & 0 deletions src/typeagent/storage/sqlite/schema.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,19 @@
);
"""

# Normalizes an ISO timestamp (any UTC offset, any fractional precision) to a
# fixed-width UTC string that sorts chronologically. Millisecond precision.
NORMALIZED_TIMESTAMP_SQL = "strftime('%Y-%m-%dT%H:%M:%f', {value})"

TIMESTAMP_INDEX_SCHEMA = """
CREATE INDEX IF NOT EXISTS idx_messages_start_timestamp ON Messages(start_timestamp);
"""

NORMALIZED_TIMESTAMP_INDEX_SCHEMA = f"""
CREATE INDEX IF NOT EXISTS idx_messages_start_timestamp_utc
ON Messages({NORMALIZED_TIMESTAMP_SQL.format(value="start_timestamp")});
"""

# Conversation metadata table (key-value pairs)
CONVERSATION_METADATA_SCHEMA = """
CREATE TABLE IF NOT EXISTS ConversationMetadata (
Expand Down Expand Up @@ -292,6 +301,7 @@ def init_db_schema(db: sqlite3.Connection) -> None:
cursor.execute(RELATED_TERMS_ALIASES_SCHEMA)
cursor.execute(RELATED_TERMS_FUZZY_SCHEMA)
cursor.execute(TIMESTAMP_INDEX_SCHEMA)
cursor.execute(NORMALIZED_TIMESTAMP_INDEX_SCHEMA)
cursor.execute(INGESTED_SOURCES_SCHEMA)
cursor.execute(CHUNK_FAILURES_SCHEMA)

Expand Down
27 changes: 17 additions & 10 deletions src/typeagent/storage/sqlite/timestampindex.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
TimestampedTextRange,
)
from ...knowpro.universal_message import format_timestamp_utc
from .schema import NORMALIZED_TIMESTAMP_SQL


class SqliteTimestampToTextRangeIndex(ITimestampToTextRangeIndex):
Expand Down Expand Up @@ -52,24 +53,26 @@ async def get_timestamp_ranges(
"""Get timestamp ranges from Messages table."""
cursor = self.db.cursor()

norm_col = NORMALIZED_TIMESTAMP_SQL.format(value="start_timestamp")
norm_arg = NORMALIZED_TIMESTAMP_SQL.format(value="?")
if end_timestamp is None:
# Single timestamp query
cursor.execute(
"""
f"""
SELECT msg_id, start_timestamp
FROM Messages
WHERE start_timestamp = ?
WHERE {norm_col} = {norm_arg}
ORDER BY msg_id
""",
(start_timestamp,),
)
else:
# Range query
# Range query (inclusive)
cursor.execute(
"""
f"""
SELECT msg_id, start_timestamp
FROM Messages
WHERE start_timestamp >= ? AND start_timestamp <= ?
WHERE {norm_col} >= {norm_arg} AND {norm_col} <= {norm_arg}
ORDER BY msg_id
""",
(start_timestamp, end_timestamp),
Expand Down Expand Up @@ -103,28 +106,32 @@ async def lookup_range(self, date_range: DateRange) -> list[TimestampedTextRange
"""Lookup messages in a date range."""
cursor = self.db.cursor()

# Convert datetime objects to ISO format strings with Z suffix
# Stored timestamps may have any UTC offset and any fractional precision,
# so raw strings don't sort chronologically. Compare via strftime(), which
# converts to UTC and renders a fixed-width value (also for existing rows).
start_timestamp = format_timestamp_utc(date_range.start)
end_timestamp = format_timestamp_utc(date_range.end) if date_range.end else None

norm_col = NORMALIZED_TIMESTAMP_SQL.format(value="start_timestamp")
norm_arg = NORMALIZED_TIMESTAMP_SQL.format(value="?")
if date_range.end is None:
# Point query
cursor.execute(
"""
f"""
SELECT msg_id, start_timestamp, chunks
FROM Messages
WHERE start_timestamp = ?
WHERE {norm_col} = {norm_arg}
ORDER BY msg_id
""",
(start_timestamp,),
)
else:
# Range query
cursor.execute(
"""
f"""
SELECT msg_id, start_timestamp, chunks
FROM Messages
WHERE start_timestamp >= ? AND start_timestamp < ?
WHERE {norm_col} >= {norm_arg} AND {norm_col} < {norm_arg}
ORDER BY msg_id
""",
(start_timestamp, end_timestamp),
Expand Down
28 changes: 28 additions & 0 deletions tests/test_email_import.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
_merge_chunks,
_split_into_paragraphs,
_text_to_chunks,
import_email_string,
)


Expand Down Expand Up @@ -100,3 +101,30 @@ def test_no_leading_separator_in_any_chunk(self) -> None:
assert not chunk.startswith(
"\n\n"
), f"chunk {chunk!r} has leading separator"


class TestEmailTimestampNormalization:
"""Email Date headers with different offsets must yield comparable timestamps."""

@staticmethod
def _timestamp(date_header: str) -> str | None:
raw = f"From: a@example.com\nTo: b@example.com\nDate: {date_header}\nSubject: s\n\nbody\n"
return import_email_string(raw, 1000).timestamp

def test_offsets_normalized_to_utc(self) -> None:
assert (
self._timestamp("Mon, 1 Jan 2024 12:00:00 +0500") == "2024-01-01T07:00:00Z"
)
assert (
self._timestamp("Mon, 1 Jan 2024 10:00:00 -0800") == "2024-01-01T18:00:00Z"
)

def test_lexicographic_order_is_chronological(self) -> None:
a = self._timestamp("Mon, 1 Jan 2024 12:00:00 +0500") # 07:00Z
b = self._timestamp("Mon, 1 Jan 2024 08:00:00 +0000") # 08:00Z
assert a is not None and b is not None and a < b

def test_unknown_zone_treated_as_utc(self) -> None:
assert (
self._timestamp("Mon, 1 Jan 2024 12:00:00 -0000") == "2024-01-01T12:00:00Z"
)
47 changes: 47 additions & 0 deletions tests/test_sqlite_indexes.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

"""Tests for SQLite index implementations with real embeddings."""

from datetime import timedelta, timezone
import os
import sqlite3
import tempfile
Expand All @@ -15,6 +16,8 @@
from typeagent.knowpro import interfaces
from typeagent.knowpro.convsettings import MessageTextIndexSettings
from typeagent.knowpro.interfaces import (
DateRange,
Datetime,
SemanticRef,
Term,
TextLocation,
Expand Down Expand Up @@ -188,6 +191,50 @@ async def test_timestamp_operations(self, sqlite_db: sqlite3.Connection):
)
assert len(results) == 2

@pytest.mark.asyncio
async def test_lookup_range_mixed_offsets_and_precision(
self, sqlite_db: sqlite3.Connection
):
"""Range lookup compares instants, not raw strings."""
index = SqliteTimestampToTextRangeIndex(sqlite_db)
stored = [
"2024-01-01T12:00:00+05:00", # 07:00:00 UTC (legacy email format)
"2024-01-01T07:00:00Z", # 07:00:00 UTC
"2024-01-01T07:00:00.250000Z",
"2024-01-01T07:00:01Z",
]
cursor = sqlite_db.cursor()
for i, ts in enumerate(stored):
cursor.execute(
"INSERT INTO Messages (msg_id, chunks, start_timestamp) VALUES (?, '[\"m\"]', ?)",
(i, ts),
)
sqlite_db.commit()

def ordinals(results):
return [r.range.start.message_ordinal for r in results]

utc = timezone.utc
dr = DateRange(
start=Datetime(2024, 1, 1, 7, tzinfo=utc),
end=Datetime(2024, 1, 1, 7, 0, 1, tzinfo=utc),
)
assert ordinals(await index.lookup_range(dr)) == [0, 1, 2]

# A fractional bound must exclude the whole-second value before it.
dr = DateRange(
start=Datetime(2024, 1, 1, 7, 0, 0, 500000, tzinfo=utc),
end=Datetime(2024, 1, 1, 8, tzinfo=utc),
)
assert ordinals(await index.lookup_range(dr)) == [3]

# Point query with a non-UTC offset matches both representations.
dr = DateRange(
start=Datetime(2024, 1, 1, 12, tzinfo=timezone(timedelta(hours=5))),
end=None,
)
assert ordinals(await index.lookup_range(dr)) == [0, 1]


class TestSqliteRelatedTermsAliases:
"""Test SqliteRelatedTermsAliases functionality."""
Expand Down
Loading
Loading