Coverage for app\routers\data.py: 29%

99 statements  

« prev     ^ index     » next       coverage.py v7.3.2, created at 2026-01-15 13:39 +0300

1from fastapi import APIRouter, Depends, HTTPException, UploadFile, File, Form 

2from fastapi.responses import FileResponse 

3from sqlalchemy.orm import Session 

4from typing import Any, List 

5import uuid 

6from datetime import datetime, UTC, timedelta 

7import os 

8import shutil 

9import tempfile 

10 

11from .. import schemas 

12from ..schemas import data as data_schemas 

13from ..db import DataExport, DataImport, User 

14from ..deps import get_db, get_current_user 

15from ..services.export_service import ExportService 

16 

17router = APIRouter() 

18 

19@router.post("/export/{export_type}", response_model=data_schemas.DataExport) 

20async def create_export( 

21 export_type: str, 

22 export_in: data_schemas.DataExportCreate, 

23 db: Session = Depends(get_db), 

24 current_user: User = Depends(get_current_user), 

25) -> Any: 

26 """ 

27 Create a new data export. 

28 """ 

29 if export_type != export_in.export_type: 

30 # Allow mismatch if URL is convenience, but better to enforce or ignore URL param? 

31 # Let's enforce consistency or override. 

32 # If user sends /export/audit_logs and body says "audit_logs", it's fine. 

33 if export_in.export_type and export_type != export_in.export_type: 

34 raise HTTPException(status_code=400, detail=f"Export type mismatch: URL={export_type}, Body={export_in.export_type}") 

35 

36 # Ensure export_type from URL is used if body is missing it (though schema requires it) 

37 # Schema DataExportBase requires export_type. 

38 

39 # Create export record 

40 db_obj = DataExport( 

41 id=str(uuid.uuid4()), 

42 export_type=export_in.export_type, 

43 format=export_in.format, 

44 filters=export_in.filters, 

45 status="pending", 

46 requested_by=current_user.id, 

47 created_at=datetime.now(UTC), 

48 file_path="", 

49 file_size=0 

50 ) 

51 db.add(db_obj) 

52 db.commit() 

53 db.refresh(db_obj) 

54 

55 # Process export (sync for now) 

56 service = ExportService(db) 

57 try: 

58 file_path = service.generate_export(db_obj) 

59 

60 # Update record 

61 db_obj.file_path = file_path 

62 db_obj.status = "completed" 

63 db_obj.completed_at = datetime.now(UTC) 

64 db_obj.file_size = os.path.getsize(file_path) 

65 db_obj.expires_at = datetime.now(UTC) + timedelta(days=7) 

66 

67 db.commit() 

68 db.refresh(db_obj) 

69 

70 return db_obj 

71 

72 except Exception as e: 

73 db_obj.status = "failed" 

74 db.commit() 

75 raise HTTPException(status_code=500, detail=str(e)) 

76 

77@router.get("/exports", response_model=List[data_schemas.DataExport]) 

78async def get_exports( 

79 skip: int = 0, 

80 limit: int = 100, 

81 db: Session = Depends(get_db), 

82 current_user: User = Depends(get_current_user), 

83) -> Any: 

84 """ 

85 Retrieve data exports. 

86 """ 

87 exports = db.query(DataExport).order_by(DataExport.created_at.desc()).offset(skip).limit(limit).all() 

88 return exports 

89 

90@router.get("/imports", response_model=List[data_schemas.DataImport]) 

91async def get_imports( 

92 skip: int = 0, 

93 limit: int = 100, 

94 db: Session = Depends(get_db), 

95 current_user: User = Depends(get_current_user), 

96) -> Any: 

97 """ 

98 Retrieve data imports. 

99 """ 

100 imports = db.query(DataImport).order_by(DataImport.created_at.desc()).offset(skip).limit(limit).all() 

101 return imports 

102 

103@router.post("/import", response_model=data_schemas.DataImport) 

104async def create_import( 

105 file: UploadFile = File(...), 

106 data_type: str = Form(...), 

107 format: str = Form(...), 

108 mode: str = Form("append"), 

109 db: Session = Depends(get_db), 

110 current_user: User = Depends(get_current_user), 

111) -> Any: 

112 """ 

113 Create a new data import. 

114 """ 

115 # Create temp file 

116 temp_dir = tempfile.gettempdir() 

117 file_ext = os.path.splitext(file.filename)[1] if file.filename else f".{format}" 

118 filename = f"import_{data_type}_{uuid.uuid4()}{file_ext}" 

119 file_path = os.path.join(temp_dir, filename) 

120 

121 with open(file_path, "wb") as buffer: 

122 shutil.copyfileobj(file.file, buffer) 

123 

124 file_size = os.path.getsize(file_path) 

125 

126 # Create import record 

127 db_obj = DataImport( 

128 id=str(uuid.uuid4()), 

129 import_type=data_type, 

130 format=format, 

131 file_path=file_path, 

132 file_size=file_size, 

133 row_count=0, # Will be updated after processing 

134 success_count=0, 

135 error_count=0, 

136 status="pending", 

137 imported_by=current_user.id, 

138 created_at=datetime.now(UTC) 

139 ) 

140 

141 db.add(db_obj) 

142 db.commit() 

143 db.refresh(db_obj) 

144 

145 # TODO: Trigger background task for processing 

146 # For now, just mark as completed to simulate success 

147 db_obj.status = "completed" 

148 db_obj.completed_at = datetime.now(UTC) 

149 db.commit() 

150 

151 return db_obj 

152 

153@router.delete("/exports/{export_id}", response_model=Any) 

154async def delete_export( 

155 export_id: str, 

156 db: Session = Depends(get_db), 

157 current_user: User = Depends(get_current_user), 

158) -> Any: 

159 """ 

160 Delete a data export. 

161 """ 

162 export = db.query(DataExport).filter(DataExport.id == export_id).first() 

163 if not export: 

164 raise HTTPException(status_code=404, detail="Export not found") 

165 

166 # Delete file if exists 

167 if export.file_path and os.path.exists(export.file_path): 

168 try: 

169 os.remove(export.file_path) 

170 except OSError: 

171 pass # Ignore file deletion errors 

172 

173 db.delete(export) 

174 db.commit() 

175 return {"status": "success"} 

176 

177@router.delete("/imports/{import_id}", response_model=Any) 

178async def delete_import( 

179 import_id: str, 

180 db: Session = Depends(get_db), 

181 current_user: User = Depends(get_current_user), 

182) -> Any: 

183 """ 

184 Delete a data import. 

185 """ 

186 import_obj = db.query(DataImport).filter(DataImport.id == import_id).first() 

187 if not import_obj: 

188 raise HTTPException(status_code=404, detail="Import not found") 

189 

190 # Delete file if exists 

191 if import_obj.file_path and os.path.exists(import_obj.file_path): 

192 try: 

193 os.remove(import_obj.file_path) 

194 except OSError: 

195 pass 

196 

197 db.delete(import_obj) 

198 db.commit() 

199 return {"status": "success"} 

200 

201@router.get("/exports/{export_id}/download") 

202async def download_export( 

203 export_id: str, 

204 db: Session = Depends(get_db), 

205 current_user: User = Depends(get_current_user), 

206) -> Any: 

207 """ 

208 Download export file. 

209 """ 

210 export = db.query(DataExport).filter(DataExport.id == export_id).first() 

211 if not export: 

212 raise HTTPException(status_code=404, detail="Export not found") 

213 

214 if not export.file_path or not os.path.exists(export.file_path): 

215 raise HTTPException(status_code=404, detail="Export file not found") 

216 

217 filename = os.path.basename(export.file_path) 

218 return FileResponse( 

219 path=export.file_path, 

220 filename=filename, 

221 media_type="application/octet-stream" 

222 ) 

223