88import tracemalloc
99import uuid
1010from functools import wraps
11- from typing import BinaryIO , cast
11+ from typing import BinaryIO , TypeVar , cast
1212
1313import typer
1414from pydantic import parse_obj_as
1818from sqlalchemy .ext .asyncio import AsyncSession
1919
2020from apps .activities .db .schemas import ActivityHistorySchema as ActivityHistory
21+ from apps .activities .db .schemas import ActivitySchema
2122from apps .activity_flows .db .schemas import ActivityFlowHistoriesSchema as FlowHistory
2223from apps .activity_flows .db .schemas import ActivityFlowItemHistorySchema as FlowItemHistory
2324from apps .activity_flows .db .schemas import ActivityFlowSchema
2425from apps .applets .db .schemas import AppletSchema
2526from apps .job .constants import JobStatus
2627from apps .job .errors import JobStatusError
2728from apps .job .service import JobService
28- from apps .schedule .db .schemas import EventSchema , FlowEventsSchema , PeriodicitySchema , UserEventsSchema
29+ from apps .schedule .db .schemas import (
30+ ActivityEventsSchema ,
31+ EventSchema ,
32+ FlowEventsSchema ,
33+ PeriodicitySchema ,
34+ UserEventsSchema ,
35+ )
2936from apps .schedule .domain .constants import PeriodicityType
3037from apps .shared .domain .base import PublicModel
3138from apps .workspaces .crud .user_applet_access import UserAppletAccessCRUD
4451PATH_PREFIX = settings .applet_ema .export_path_prefix
4552PATH_FLOW_FILE_NAME = settings .applet_ema .export_flow_file_name
4653PATH_USER_FLOW_SCHEDULE_FILE_NAME = settings .applet_ema .export_user_flow_schedule_file_name
54+ PATH_USER_ACTIVITY_SCHEDULE_FILE_NAME = settings .applet_ema .export_user_activity_schedule_file_name
4755
4856
4957# Not ISO
5563OUTPUT_TIME_FORMAT = "%H:%M"
5664
5765
58- class OutputRow (PublicModel ):
66+ class FlowEventOutputRow (PublicModel ):
5967 applet_id : uuid .UUID
6068 date_prior_day : datetime .date
6169 user_id : uuid .UUID
70+ secret_user_id : str | uuid .UUID
6271 flow_id : uuid .UUID
6372 flow_name : str
6473 applet_version : str
@@ -69,12 +78,26 @@ class OutputRow(PublicModel):
6978 event_id : uuid .UUID
7079
7180
81+ class ActivityEventOutputRow (PublicModel ):
82+ applet_id : uuid .UUID
83+ date_prior_day : datetime .date
84+ user_id : uuid .UUID
85+ secret_user_id : str | uuid .UUID
86+ activity_id : uuid .UUID
87+ activity_name : str
88+ applet_version : str
89+ scheduled_date : datetime .date
90+ schedule_start_time : str
91+ schedule_end_time : str
92+ # TODO: remove after debug
93+ event_id : uuid .UUID
94+
95+
7296class RawRow (PublicModel ):
7397 applet_id : uuid .UUID
7498 date : datetime .date
7599 user_id : uuid .UUID
76- flow_id : uuid .UUID
77- flow_name : str
100+ secret_user_id : str | uuid .UUID
78101 applet_version : str
79102 schedule_start_time : datetime .time
80103 schedule_end_time : datetime .time
@@ -89,6 +112,19 @@ def is_crossday_event(self) -> bool:
89112 return self .schedule_start_time > self .schedule_end_time
90113
91114
115+ TRawRow = TypeVar ("TRawRow" , bound = RawRow )
116+
117+
118+ class FlowEventRawRow (RawRow ):
119+ flow_id : uuid .UUID
120+ flow_name : str
121+
122+
123+ class ActivityEventRawRow (RawRow ):
124+ activity_id : uuid .UUID
125+ activity_name : str
126+
127+
92128def is_last_day_of_month (date : datetime .date ):
93129 mdays = calendar .mdays .copy () # type: ignore[attr-defined]
94130 if calendar .isleap (date .year ):
@@ -198,7 +234,7 @@ async def export_flows():
198234
199235
200236##### Daily user flow schedule stuff
201- async def get_user_flow_events (session : AsyncSession , scheduled_date : datetime .date ) -> list [RawRow ]:
237+ async def get_user_flow_events (session : AsyncSession , scheduled_date : datetime .date ) -> list [FlowEventRawRow ]:
202238 cte = (
203239 select (
204240 EventSchema .applet_id ,
@@ -231,6 +267,7 @@ async def get_user_flow_events(session: AsyncSession, scheduled_date: datetime.d
231267 AppletSchema .id .label ("applet_id" ),
232268 literal (scheduled_date , Date ).label ("date" ),
233269 UserAppletAccessSchema .user_id .label ("user_id" ),
270+ UserAppletAccessSchema .respondent_secret_id .label ("secret_user_id" ), # type: ignore[attr-defined] # noqa: E501
234271 ActivityFlowSchema .id .label ("flow_id" ),
235272 ActivityFlowSchema .name .label ("flow_name" ),
236273 AppletSchema .version .label ("applet_version" ),
@@ -268,11 +305,11 @@ async def get_user_flow_events(session: AsyncSession, scheduled_date: datetime.d
268305 )
269306 db_result = await session .execute (query )
270307 result = db_result .mappings ().all ()
271- return parse_obj_as (list [RawRow ], result )
308+ return parse_obj_as (list [FlowEventRawRow ], result )
272309
273310
274- def filter_events (raw_events_rows : list [RawRow ], schedule_date : datetime .date ) -> list [RawRow ]: # noqa: C901
275- filtered : list [RawRow ] = []
311+ def filter_events (raw_events_rows : list [TRawRow ], schedule_date : datetime .date ) -> list [TRawRow ]: # noqa: C901
312+ filtered : list [TRawRow ] = []
276313 for row in raw_events_rows :
277314 # TODO: patch events with periodicity WEEKDAYS, WEEKLY, some events don't have start_date and end_date
278315 # (the issue is in migrated data).
@@ -293,7 +330,11 @@ def filter_events(raw_events_rows: list[RawRow], schedule_date: datetime.date) -
293330 if schedule_date >= schedule_start_date and schedule_date <= row .end_date :
294331 filtered .append (row )
295332 case PeriodicityType .WEEKDAYS :
296- last_weekday = FRIDAY_WEEKDAY if not row .is_crossday_event else SATURDAY_WEEKDAY
333+ last_weekday = FRIDAY_WEEKDAY
334+ if row .is_crossday_event :
335+ last_weekday = SATURDAY_WEEKDAY
336+ if row .end_date .weekday () == FRIDAY_WEEKDAY :
337+ row .end_date += datetime .timedelta (days = 1 )
297338 if (
298339 schedule_date .weekday () <= last_weekday
299340 and schedule_date >= row .start_date
@@ -358,7 +399,7 @@ async def export_flow_schedule(
358399 ),
359400):
360401 ensure_configured ()
361- scheduled_date = run_date .date () if run_date else datetime .date .today () - datetime . timedelta ( days = 1 )
402+ scheduled_date = run_date .date () if run_date else datetime .date .today ()
362403
363404 job_name = f"export_flow_schedule_{ scheduled_date } "
364405
@@ -392,10 +433,11 @@ async def export_flow_schedule(
392433 print (f"Num filtered rows is { len (filtered )} " )
393434 result = []
394435 for row in filtered :
395- outrow = OutputRow (
436+ outrow = FlowEventOutputRow (
396437 applet_id = row .applet_id ,
397438 date_prior_day = scheduled_date ,
398439 user_id = row .user_id ,
440+ secret_user_id = row .secret_user_id ,
399441 flow_id = row .flow_id ,
400442 flow_name = row .flow_name ,
401443 applet_version = row .applet_version ,
@@ -407,7 +449,7 @@ async def export_flow_schedule(
407449 result .append (outrow )
408450
409451 cdn_client = await get_operations_bucket ()
410- unique_prefix = f"{ APPLET_ID } /flow_schedule "
452+ unique_prefix = f"{ APPLET_ID } /flow-schedule "
411453
412454 prev_filename = PATH_USER_FLOW_SCHEDULE_FILE_NAME .format (date = scheduled_date - datetime .timedelta (days = 1 ))
413455 prev_key = cdn_client .generate_key (PATH_PREFIX , unique_prefix , prev_filename )
@@ -425,6 +467,7 @@ async def export_flow_schedule(
425467 f .seek (0 , io .SEEK_END )
426468 create_csv (result , append_to = f )
427469 with open (path , "rb" ) as f :
470+ print (f"Upload file to the { key } " )
428471 await cdn_client .upload (key , f )
429472
430473 os .remove (path )
@@ -442,3 +485,177 @@ async def export_flow_schedule(
442485 tracemalloc .stop ()
443486 print ("Flow schedule export finished" )
444487 print ("Peak memory usage:" , peak )
488+
489+
490+ ##### Daily user activity schedule stuff
491+ async def get_user_activity_events (session : AsyncSession , scheduled_date : datetime .date ) -> list [ActivityEventRawRow ]:
492+ cte = (
493+ select (
494+ EventSchema .applet_id ,
495+ EventSchema .id .label ("event_id" ),
496+ UserEventsSchema .user_id ,
497+ ActivityEventsSchema .activity_id ,
498+ PeriodicitySchema .type .label ("event_type" ),
499+ case (
500+ (
501+ PeriodicitySchema .type .in_ (("WEEKDAYS" , "DAILY" )),
502+ scheduled_date ,
503+ ),
504+ (PeriodicitySchema .type .in_ (("WEEKLY" , "MONTHLY" )), PeriodicitySchema .start_date ),
505+ else_ = PeriodicitySchema .selected_date ,
506+ ).label ("selected_date" ),
507+ PeriodicitySchema .start_date ,
508+ PeriodicitySchema .end_date ,
509+ EventSchema .start_time ,
510+ EventSchema .end_time ,
511+ )
512+ .select_from (EventSchema )
513+ .join (UserEventsSchema , UserEventsSchema .event_id == EventSchema .id )
514+ .join (PeriodicitySchema , PeriodicitySchema .id == EventSchema .periodicity_id )
515+ .join (ActivityEventsSchema , ActivityEventsSchema .event_id == EventSchema .id )
516+ .where (EventSchema .is_deleted == false (), PeriodicitySchema .type != PeriodicityType .ALWAYS )
517+ ).cte ("user_activity_events" )
518+
519+ query = (
520+ select (
521+ AppletSchema .id .label ("applet_id" ),
522+ literal (scheduled_date , Date ).label ("date" ),
523+ UserAppletAccessSchema .user_id .label ("user_id" ),
524+ UserAppletAccessSchema .respondent_secret_id .label ("secret_user_id" ), # type: ignore[attr-defined] # noqa: E501
525+ ActivitySchema .id .label ("activity_id" ),
526+ ActivitySchema .name .label ("activity_name" ),
527+ AppletSchema .version .label ("applet_version" ),
528+ cte .c .event_id .label ("event_id" ),
529+ cte .c .event_type .label ("event_type" ),
530+ cte .c .selected_date .label ("selected_date" ),
531+ cte .c .start_date .label ("start_date" ),
532+ cte .c .end_date .label ("end_date" ),
533+ cte .c .start_time .label ("schedule_start_time" ),
534+ cte .c .end_time .label ("schedule_end_time" ),
535+ )
536+ .select_from (AppletSchema )
537+ .join (
538+ UserAppletAccessSchema ,
539+ and_ (
540+ UserAppletAccessSchema .applet_id == AppletSchema .id ,
541+ UserAppletAccessSchema .role == Role .RESPONDENT ,
542+ ),
543+ )
544+ .join (ActivitySchema , ActivitySchema .applet_id == AppletSchema .id )
545+ .outerjoin (
546+ cte ,
547+ and_ (
548+ cte .c .applet_id == AppletSchema .id ,
549+ cte .c .user_id == UserAppletAccessSchema .user_id ,
550+ cte .c .activity_id == ActivitySchema .id ,
551+ ),
552+ )
553+ .where (
554+ AppletSchema .id == uuid .UUID (APPLET_ID ),
555+ cte .c .event_id != null (),
556+ ActivitySchema .is_hidden == false (),
557+ )
558+ .order_by (UserAppletAccessSchema .user_id , ActivitySchema .name )
559+ )
560+ db_result = await session .execute (query )
561+ result = db_result .mappings ().all ()
562+ return parse_obj_as (list [ActivityEventRawRow ], result )
563+
564+
565+ @app .command (short_help = "Export daily user activity schedule events to csv" )
566+ @coro
567+ async def export_activity_schedule (
568+ run_date : datetime .datetime = typer .Argument (None , help = "run date" ),
569+ force : bool = typer .Option (
570+ False ,
571+ "--force" ,
572+ "-f" ,
573+ help = "Force run even if job executed before" ,
574+ ),
575+ ):
576+ ensure_configured ()
577+ scheduled_date = run_date .date () if run_date else datetime .date .today ()
578+
579+ job_name = f"export_activity_schedule_{ scheduled_date } "
580+
581+ session_maker = session_manager .get_session ()
582+ async with session_maker () as session :
583+ owner_role = await UserAppletAccessCRUD (session ).get_applet_owner (uuid .UUID (APPLET_ID ))
584+ owner_id = owner_role .user_id
585+
586+ job_service = JobService (session , owner_id )
587+ async with atomic (session ):
588+ try :
589+ job = await job_service .get_or_create_owned (
590+ job_name , JobStatus .in_progress , accept_statuses = [JobStatus .error ]
591+ )
592+ except JobStatusError as e :
593+ # prevent task execution if it's in progress or executed previously
594+ job = e .job
595+ if job .status == JobStatus .in_progress or not force :
596+ raise
597+ if job .status != JobStatus .in_progress :
598+ await job_service .change_status (job .id , JobStatus .in_progress )
599+ print ("Activity schedule export start" )
600+ tracemalloc .start ()
601+
602+ try :
603+ session_maker = session_manager .get_session ()
604+ async with session_maker () as session :
605+ raw_data = await get_user_activity_events (session , scheduled_date )
606+ print (f"Num raw rows is { len (raw_data )} " )
607+ filtered = filter_events (raw_data , scheduled_date )
608+ print (f"Num filtered rows is { len (filtered )} " )
609+ result = []
610+ for row in filtered :
611+ outrow = ActivityEventOutputRow (
612+ applet_id = row .applet_id ,
613+ date_prior_day = scheduled_date ,
614+ user_id = row .user_id ,
615+ secret_user_id = row .secret_user_id ,
616+ activity_id = row .activity_id ,
617+ activity_name = row .activity_name ,
618+ applet_version = row .applet_version ,
619+ scheduled_date = scheduled_date ,
620+ schedule_start_time = row .schedule_start_time .strftime (OUTPUT_TIME_FORMAT ),
621+ schedule_end_time = row .schedule_end_time .strftime (OUTPUT_TIME_FORMAT ),
622+ event_id = row .event_id ,
623+ ).dict ()
624+ result .append (outrow )
625+
626+ cdn_client = await get_operations_bucket ()
627+ unique_prefix = f"{ APPLET_ID } /activity-schedule"
628+
629+ prev_filename = PATH_USER_ACTIVITY_SCHEDULE_FILE_NAME .format (date = scheduled_date - datetime .timedelta (days = 1 ))
630+ prev_key = cdn_client .generate_key (PATH_PREFIX , unique_prefix , prev_filename )
631+
632+ filename = PATH_USER_ACTIVITY_SCHEDULE_FILE_NAME .format (date = scheduled_date )
633+ key = cdn_client .generate_key (PATH_PREFIX , unique_prefix , filename )
634+
635+ path = settings .uploads_dir / filename
636+
637+ with open (path , "wb" ) as f :
638+ try :
639+ cdn_client .download (prev_key , f )
640+ except ObjectNotFoundError :
641+ pass
642+ f .seek (0 , io .SEEK_END )
643+ create_csv (result , append_to = f )
644+ with open (path , "rb" ) as f :
645+ print (f"Upload file to the { key } " )
646+ await cdn_client .upload (key , f )
647+
648+ os .remove (path )
649+ async with session_maker () as session :
650+ async with atomic (session ):
651+ await JobService (session , owner_id ).change_status (job .id , JobStatus .success )
652+ except Exception as e :
653+ async with session_maker () as session :
654+ async with atomic (session ):
655+ await JobService (session , owner_id ).change_status (job .id , JobStatus .error , {"error" : str (e )})
656+ raise
657+
658+ _ , peak = tracemalloc .get_traced_memory ()
659+ tracemalloc .stop ()
660+ print ("Activity schedule export finished" )
661+ print ("Peak memory usage:" , peak )
0 commit comments