1414ACTIVITY_PAGE_MAX = 200
1515ACTIVITY_SUMMARY_MAX = 5000
1616ACTIVITY_EXPORT_MAX = 10000
17+ ACTIVITY_EXPORT_HEADER = (
18+ "timestamp" , "id" , "user_id" , "activity_type" , "workspace_type" , "user_name" , "user_email" ,
19+ "activity" , "summary" , "workspace_id" , "workspace_name" , "raw_json" ,
20+ )
21+ ACTIVITY_EXPORT_READABLE_COLUMNS = 6
1722ACTIVITY_ORDER = " ORDER BY c.timestamp DESC, c.id DESC, c.user_id DESC"
1823ACTIVITY_COMPOSITE_INDEX = [
1924 {"path" : "/timestamp" , "order" : "descending" },
2833 "group.group_name" , "workspace_context.group_id" ,
2934 "workspace_context.public_workspace_id" ,
3035)
36+ # Cosmos SQL keywords that cannot be used as dotted property names; the query grammar only
37+ # accepts ALL, FIRST and LAST there. A dotted reserved word (c.group.group_id) is a syntax
38+ # error that rejects the whole query, so those segments are written as c['group'].
39+ COSMOS_RESERVED_WORDS = frozenset ({
40+ "AND" , "ARRAY" , "AS" , "ASC" , "BETWEEN" , "BY" , "DESC" , "DISTINCT" , "ESCAPE" , "EXISTS" ,
41+ "FALSE" , "FROM" , "GROUP" , "IN" , "JOIN" , "LEFT" , "LIKE" , "LIMIT" , "NOT" , "NULL" , "OFFSET" ,
42+ "OR" , "ORDER" , "RANK" , "RIGHT" , "SELECT" , "TOP" , "TRUE" , "UDF" , "UNDEFINED" , "VALUE" , "WHERE" ,
43+ })
44+
45+
46+ def cosmos_property_path (path , root = "c" ):
47+ """Return a Cosmos property reference that is valid even when a segment is a keyword."""
48+ expression = root
49+ for segment in path .split ("." ):
50+ expression += f"['{ segment } ']" if segment .upper () in COSMOS_RESERVED_WORDS else f".{ segment } "
51+ return expression
52+
53+
54+ NESTED_GROUP_ID = cosmos_property_path ("group.group_id" )
55+ # Where the writers record who acted. Records from approvals, membership and status changes
56+ # often have no top-level user_id, so the person filter and the people search match all of
57+ # them. The display module resolves the table's Person column from the same fields.
58+ ACTIVITY_ACTOR_FIELDS = (
59+ "user_id" , "admin_user_id" , "requester_id" , "added_by_user_id" , "changed_by_user_id" ,
60+ "changed_by.user_id" , "removed_by.user_id" , "admin.user_id" , "actor.user_id" ,
61+ )
62+ # Where the writers record a group or public workspace. Status changes and member removals
63+ # nest it (group.group_id, public_workspace.*), approvals and membership audits store it at
64+ # the top level, user agreements store workspace_context.<type>_workspace_id, and public
65+ # workspace ownership approvals store only a bare workspace_id. The display module resolves
66+ # the Workspace column from the same locations.
67+ GROUP_REFERENCE_FIELDS = (
68+ "workspace_context.group_id" , "group_id" , "group.group_id" , "workspace_context.group_workspace_id" ,
69+ )
70+ PUBLIC_REFERENCE_FIELDS = (
71+ "workspace_context.public_workspace_id" , "public_workspace_id" , "public_workspace.public_workspace_id" ,
72+ "public_workspace.workspace_id" , "workspace_id" ,
73+ )
3174
3275
3376def parse_activity_filters (args , now = None ):
@@ -60,7 +103,14 @@ def parse_activity_filters(args, now=None):
60103 return result
61104
62105
63- def activity_query_context (filters ):
106+ def _either (fields , template ):
107+ """OR the same predicate over several property paths, bracket-quoting reserved words."""
108+ return " OR " .join (template .format (path = cosmos_property_path (field )) for field in fields )
109+
110+
111+ def activity_query_context (filters , search_user_ids = ()):
112+ """Build the WHERE clause. search_user_ids are people whose name or email matched the
113+ search; they widen the search only and stay outside the cursor's filter scope."""
64114 end_exclusive = (date .fromisoformat (filters ["end_date" ]) + timedelta (days = 1 )).isoformat ()
65115 clauses = ["IS_STRING(c.timestamp)" , "IS_STRING(c.id)" ,
66116 "(IS_STRING(c.user_id) OR IS_NULL(c.user_id) OR NOT IS_DEFINED(c.user_id))" ,
@@ -74,21 +124,22 @@ def add(clause, name, value):
74124 parameters .append ({"name" : name , "value" : value })
75125
76126 add ("ARRAY_CONTAINS(@types, c.activity_type)" , "@types" , filters ["activity_types" ])
77- add ("(c.user_id = @user OR c.changed_by.user_id = @user OR c.admin_user_id = @user)" ,
78- "@user" , filters ["user_id" ])
127+ add (f"({ _either (ACTIVITY_ACTOR_FIELDS , '{path} = @user' )} )" , "@user" , filters ["user_id" ])
79128 workspace_type = filters ["workspace_type" ]
80- if workspace_type == "public" :
81- add ("c.workspace_type IN ('public', 'public_workspace')" , "@workspace_type" , workspace_type )
82- else :
83- add ("c.workspace_type = @workspace_type" , "@workspace_type" , workspace_type )
84129 group_id = filters ["group_id" ] or (filters ["workspace_id" ] if workspace_type == "group" else "" )
85130 public_id = filters ["public_workspace_id" ] or (filters ["workspace_id" ] if workspace_type == "public" else "" )
86- add ("(c.workspace_context.group_id = @group OR c.group_id = @group OR c.group.group_id = @group)" ,
87- "@group" , group_id )
88- add ("(c.workspace_context.public_workspace_id = @public OR c.public_workspace_id = @public)" ,
89- "@public" , public_id )
131+ # A specific group or public workspace matches every record that references it, however
132+ # its writer recorded the workspace type; a type on its own matches the same references.
90133 if workspace_type == "personal" :
134+ add ("c.workspace_type = @workspace_type" , "@workspace_type" , workspace_type )
91135 add ("c.user_id = @personal" , "@personal" , filters ["workspace_id" ])
136+ elif workspace_type == "group" and not group_id :
137+ clauses .append (f"(c.workspace_type = 'group' OR { _either (GROUP_REFERENCE_FIELDS , 'IS_STRING({path})' )} )" )
138+ elif workspace_type == "public" and not public_id :
139+ clauses .append ("(c.workspace_type IN ('public', 'public_workspace') OR "
140+ f"{ _either (PUBLIC_REFERENCE_FIELDS , 'IS_STRING({path})' )} )" )
141+ add (f"({ _either (GROUP_REFERENCE_FIELDS , '{path} = @group' )} )" , "@group" , group_id )
142+ add (f"({ _either (PUBLIC_REFERENCE_FIELDS , '{path} = @public' )} )" , "@public" , public_id )
92143 add ("c.token_type = @token_type" , "@token_type" , filters ["token_type" ])
93144 add ("c.usage.model = @model" , "@model" , filters ["model" ])
94145 if filters ["status" ] == "failed" :
@@ -99,7 +150,10 @@ def add(clause, name, value):
99150 add ("(c.status = @status OR c.status_change.new_status = @status OR c.document.status = @status)" ,
100151 "@status" , filters ["status" ])
101152 if filters ["search" ]:
102- search = " OR " .join (f"CONTAINS(c.{ field } , @search, true)" for field in ACTIVITY_SEARCH_FIELDS )
153+ search = _either (ACTIVITY_SEARCH_FIELDS , "CONTAINS({path}, @search, true)" )
154+ if search_user_ids :
155+ search += " OR " + _either (ACTIVITY_ACTOR_FIELDS , "ARRAY_CONTAINS(@search_people, {path})" )
156+ parameters .append ({"name" : "@search_people" , "value" : list (search_user_ids )})
103157 add (f"({ search } )" , "@search" , filters ["search" ])
104158 return " AND " .join (clauses ), parameters
105159
@@ -138,9 +192,9 @@ def decode_activity_cursor(value, filters):
138192 raise ValueError ("Invalid activity cursor" ) from ex
139193
140194
141- def query_activity_rows (container , filters , limit , cursor = None , snapshot = None , projection = "*" ):
195+ def query_activity_rows (container , filters , limit , cursor = None , snapshot = None , projection = "*" , search_user_ids = () ):
142196 """Three-part key: IDs are only unique inside a user partition, not globally."""
143- where , parameters = activity_query_context (filters )
197+ where , parameters = activity_query_context (filters , search_user_ids )
144198 snapshot = cursor ["snapshot" ] if cursor else snapshot or datetime .now (timezone .utc ).isoformat ()
145199 # Compare calendar instants from both legacy naive-UTC and aware-UTC writers.
146200 # The cursor keeps the stored timestamp verbatim so lexical ordering and seeking agree.
@@ -170,11 +224,13 @@ def query_activity_rows(container, filters, limit, cursor=None, snapshot=None, p
170224 return rows , snapshot
171225
172226
173- def activity_page (container , filters , page_size = 50 , cursor_value = None ):
227+ def activity_page (container , filters , page_size = 50 , cursor_value = None , search_user_ids = () ):
174228 if not 1 <= page_size <= ACTIVITY_PAGE_MAX :
175229 raise ValueError ("Invalid page size" )
176230 cursor = decode_activity_cursor (cursor_value , filters ) if cursor_value else None
177- rows , snapshot = query_activity_rows (container , filters , page_size + 1 , cursor = cursor )
231+ rows , snapshot = query_activity_rows (
232+ container , filters , page_size + 1 , cursor = cursor , search_user_ids = search_user_ids ,
233+ )
178234 more = len (rows ) > page_size
179235 records = rows [:page_size ]
180236 return {
@@ -183,10 +239,11 @@ def activity_page(container, filters, page_size=50, cursor_value=None):
183239 }
184240
185241
186- def activity_summary (container , filters ):
242+ def activity_summary (container , filters , search_user_ids = () ):
187243 """Bounded projection, not COUNT/GROUP BY scans. Disclose sampling in the contract."""
188244 rows , snapshot = query_activity_rows (
189245 container , filters , ACTIVITY_SUMMARY_MAX + 1 , projection = "c.timestamp, c.id, c.user_id, c.activity_type" ,
246+ search_user_ids = search_user_ids ,
190247 )
191248 truncated = len (rows ) > ACTIVITY_SUMMARY_MAX
192249 rows = rows [:ACTIVITY_SUMMARY_MAX ]
@@ -216,27 +273,41 @@ def activity_csv_cell(value):
216273 return text
217274
218275
219- def activity_csv_stream (container , filters , first_rows , snapshot ):
220- """Stream page-sized reads with a final status row when the export hits its cap."""
276+ def activity_csv_stream (container , filters , first_rows , snapshot , search_user_ids = (), present_rows = None ):
277+ """Stream page-sized reads with a final status row when the export hits its cap.
278+
279+ The first five columns and the trailing raw JSON keep their 0.261.284 positions. The
280+ readable columns between them come from present_rows(rows), which returns one
281+ (user_name, user_email, activity, summary, workspace_id, workspace_name) per row.
282+ """
221283 output = StringIO ()
222284 writer = csv .writer (output )
285+ blank = ("" ,) * ACTIVITY_EXPORT_READABLE_COLUMNS
223286
224287 def line (values ):
225288 output .seek (0 )
226289 output .truncate (0 )
227290 writer .writerow ([activity_csv_cell (value ) for value in values ])
228291 return output .getvalue ()
229292
230- yield line (( "timestamp" , "id" , "user_id" , "activity_type" , "workspace_type" , "raw_json" ) )
293+ yield line (ACTIVITY_EXPORT_HEADER )
231294 rows = first_rows
232295 count = 0
233296 while rows :
234- for row in rows [:ACTIVITY_EXPORT_MAX - count ]:
297+ batch = rows [:ACTIVITY_EXPORT_MAX - count ]
298+ readable = list (present_rows (batch )) if present_rows else []
299+ if len (readable ) != len (batch ):
300+ readable = [blank ] * len (batch )
301+ for row , columns in zip (batch , readable ):
235302 yield line ((row .get ("timestamp" ), row .get ("id" ), row .get ("user_id" ), row .get ("activity_type" ),
236- row .get ("workspace_type" ), json .dumps (row , ensure_ascii = False )))
303+ row .get ("workspace_type" ), * columns , json .dumps (row , ensure_ascii = False )))
237304 count += 1
238305 if count >= ACTIVITY_EXPORT_MAX :
239- yield line (("" , "" , "" , "export_limit_reached" , "" , f"Export capped at { ACTIVITY_EXPORT_MAX } rows; narrow the filters." ))
306+ yield line (("" , "" , "" , "export_limit_reached" , "" , * blank ,
307+ f"Export capped at { ACTIVITY_EXPORT_MAX } rows; narrow the filters." ))
240308 return
241309 cursor = decode_activity_cursor (encode_activity_cursor (rows [- 1 ], filters , snapshot ), filters )
242- rows , _ = query_activity_rows (container , filters , min (ACTIVITY_PAGE_MAX , ACTIVITY_EXPORT_MAX - count ), cursor )
310+ rows , _ = query_activity_rows (
311+ container , filters , min (ACTIVITY_PAGE_MAX , ACTIVITY_EXPORT_MAX - count ), cursor ,
312+ search_user_ids = search_user_ids ,
313+ )
0 commit comments